"""The live lane: a loopback bridge into the running Cinema 4D.

`three-blocks c4d graph|summary` (no file argument) reads the OPEN document
through this — the C4D analog of blender_mcp's :9876 addon socket, opt-in via
TB_C4D_BRIDGE=1 at launch (the sandbox sets it).

Security posture: loopback only, and every request must carry the random token
this process wrote to `~/.three-blocks/c4d/bridge.json` (0600) — same trust
boundary as the credential store itself. No arbitrary-code op exists.

Threading: the socket thread only parses requests; each op runs on the MAIN
thread — the socket thread queues a job + `SpecialEventAdd`, the MessageData
handler in three_blocks.pyp calls `drain_jobs()`, and the socket thread waits
on a per-job reply queue. Scene access off-main is forbidden in C4D.
"""

from __future__ import annotations

import json
import os
import queue
import secrets
import socket
import threading

import c4d
import c4d.documents

from . import events, host, prefs, scene_report
from ._core import config

_JOBS: "queue.Queue[tuple]" = queue.Queue()
_server: socket.socket | None = None

REQUEST_LIMIT = 64 * 1024
JOB_TIMEOUT_SECONDS = 15.0


def handshake_path():
    return config.config_dir() / "c4d" / "bridge.json"


def enabled() -> bool:
    """Auto-start: TB_C4D_BRIDGE=1 (sandbox/CI) or the remembered panel toggle."""
    if os.environ.get("TB_C4D_BRIDGE") == "1":
        return True
    try:
        return prefs.bridge_autostart()
    except Exception:
        return False


def running() -> bool:
    return _server is not None


def port() -> int | None:
    return _server.getsockname()[1] if _server is not None else None


# --- ops (main thread) --------------------------------------------------------


def _active_document():
    doc = c4d.documents.GetActiveDocument()
    if doc is None:
        raise RuntimeError("No active Cinema 4D document.")
    return doc


def _op_ping(params: dict) -> dict:
    doc = c4d.documents.GetActiveDocument()
    return {
        "label": host.label(),
        "hub": host.hub_version(),
        "document": doc.GetDocumentName() if doc else None,
        "path": doc.GetDocumentPath() if doc else None,
        "dirty": bool(doc.GetChanged()) if doc else False,
    }


# Host-level ops only. Scene ops live in scene_report.OPS (register_op) so new
# ops never touch this socket code — the c4dpy lane dispatches the same table.
_OPS = {"ping": _op_ping}


def drain_jobs() -> None:
    """Run queued ops on the main thread (called from MessageData.CoreMessage)."""
    while True:
        try:
            op, params, reply = _JOBS.get_nowait()
        except queue.Empty:
            return
        try:
            handler = _OPS.get(op)
            if handler is not None:
                reply.put({"ok": True, "result": handler(params or {})})
                continue
            scene_op = scene_report.OPS.get(op)
            if scene_op is None:
                reply.put({"ok": False, "error": f"Unknown op '{op}'."})
            else:
                reply.put({"ok": True, "result": scene_op(_active_document(), params or {})})
        except Exception as error:
            reply.put({"ok": False, "error": str(error)})


# --- socket plumbing (worker threads) ----------------------------------------


def _serve_connection(connection: socket.socket, token: str) -> None:
    with connection:
        connection.settimeout(JOB_TIMEOUT_SECONDS)
        try:
            raw = b""
            while not raw.endswith(b"\n") and len(raw) < REQUEST_LIMIT:
                chunk = connection.recv(4096)
                if not chunk:
                    break
                raw += chunk
            request = json.loads(raw.decode("utf-8"))
            if not secrets.compare_digest(str(request.get("token") or ""), token):
                connection.sendall(b'{"ok": false, "error": "Bad bridge token."}\n')
                return
            reply: "queue.Queue[dict]" = queue.Queue()
            _JOBS.put((request.get("op"), request.get("params"), reply))
            c4d.SpecialEventAdd(prefs.HUB_ID)
            try:
                response = reply.get(timeout=JOB_TIMEOUT_SECONDS)
            except queue.Empty:
                response = {"ok": False, "error": "Cinema 4D did not answer in time (busy?)."}
            connection.sendall(json.dumps(response).encode("utf-8") + b"\n")
        except Exception:
            try:
                connection.sendall(b'{"ok": false, "error": "Bad bridge request."}\n')
            except Exception:
                pass


def start() -> int | None:
    """Bind the loopback server, write the handshake file, serve forever."""
    global _server
    if _server is not None:
        return _server.getsockname()[1]
    token = secrets.token_hex(16)
    server = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
    server.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
    server.bind(("127.0.0.1", int(os.environ.get("TB_C4D_BRIDGE_PORT") or 0)))
    server.listen(4)
    port = server.getsockname()[1]
    path = handshake_path()
    path.parent.mkdir(parents=True, exist_ok=True)
    path.write_text(
        json.dumps({"port": port, "token": token, "pid": os.getpid(), "label": host.label()}) + "\n",
        encoding="utf-8",
    )
    try:
        os.chmod(path, 0o600)
    except OSError:
        pass
    _server = server

    def accept_loop():
        while True:
            try:
                connection, _ = server.accept()
            except OSError:
                return  # closed on stop()
            threading.Thread(target=_serve_connection, args=(connection, token), daemon=True).start()

    events.spawn(accept_loop)
    print(f"three_blocks: bridge listening on 127.0.0.1:{port}")
    return port


def stop() -> None:
    global _server
    if _server is not None:
        try:
            _server.close()
        finally:
            _server = None
    handshake_path().unlink(missing_ok=True)
